feat(kafka-connect): reset connector using REST API, remove custom resetter - #646
Conversation
There was a problem hiding this comment.
Pull request overview
This PR updates the Kafka Connect connector reset/clean flow to use the Kafka Connect REST API for offset resets (per KIP-875 / Kafka 3.6+) and removes the custom Helm-based “resetter” mechanism along with its configuration surface area and documentation/testing artifacts.
Changes:
- Implement connector offset reset via Kafka Connect REST API (
KafkaConnectHandler.reset_connector+ConnectWrapper.reset_offset). - Remove the Helm resetter model/config (
resetter_namespace,resetter_values) from connector components, schemas, docs, and snapshot expectations. - Update unit tests to reflect the new reset behavior and logging patterns.
Reviewed changes
Copilot reviewed 30 out of 30 changed files in this pull request and generated 3 comments.
Show a summary per file
| File | Description |
|---|---|
| tests/pipeline/test_generate.py | Removes tests that validated resetter-related value substitution. |
| tests/pipeline/snapshots/test_generate/test_with_env_defaults/pipeline.yaml | Updates snapshot output to omit resetter_values. |
| tests/pipeline/snapshots/test_generate/test_read_from_component/pipeline.yaml | Updates snapshot output to omit resetter_values in multiple places. |
| tests/pipeline/snapshots/test_generate/test_kafka_connect_sink_weave_from_topics/pipeline.yaml | Updates snapshot output to omit resetter_values. |
| tests/pipeline/snapshots/test_generate/test_inflate_pipeline/pipeline.yaml | Updates snapshot output to omit resetter_values. |
| tests/pipeline/snapshots/test_example/test_generate/word-count/pipeline.yaml | Updates snapshot output to omit resetter_values. |
| tests/pipeline/snapshots/test_example/test_generate/atm-fraud/pipeline.yaml | Updates snapshot output to omit resetter_values. |
| tests/pipeline/resources/resetter_values/pipeline.yaml | Removes resetter-values test resource file. |
| tests/pipeline/resources/resetter_values/pipeline_connector_only.yaml | Removes resetter-values test resource file. |
| tests/pipeline/resources/resetter_values/defaults.yaml | Removes resetter-values defaults test resource file. |
| tests/components/test_kafka_source_connector.py | Removes resetter-specific assertions and updates reset/clean expectations to use reset_connector. |
| tests/components/test_kafka_sink_connector.py | Removes resetter-specific assertions and updates reset/clean expectations to use reset_connector. |
| tests/components/test_kafka_connector.py | Removes resetter_namespace usage from base connector tests. |
| tests/component_handlers/kafka_connect/test_connect_wrapper.py | Wraps calls in caplog.at_level to assert INFO logs reliably. |
| tests/component_handlers/kafka_connect/test_connect_handler.py | Adds coverage for KafkaConnectHandler.reset_connector (dry-run, missing connector, concurrent disappearance). |
| kpops/components/base_components/kafka_connector.py | Removes resetter Helm implementation and routes reset/clean through connector_handler.reset_connector. |
| kpops/component_handlers/kafka_connect/kafka_connect_handler.py | Adds reset_connector() implementation using REST operations (stop + offsets reset) with dry-run support. |
| kpops/component_handlers/kafka_connect/connect_wrapper.py | Adds reset_offset() REST call wrapper. |
| docs/docs/schema/pipeline.json | Removes resetter_namespace / resetter_values from published pipeline schema. |
| docs/docs/schema/defaults.json | Removes resetter_namespace / resetter_values from defaults schema. |
| docs/docs/resources/pipeline-defaults/defaults.yaml | Removes resetter-related defaults examples. |
| docs/docs/resources/pipeline-defaults/defaults-kafka-connector.yaml | Removes resetter-related defaults examples. |
| docs/docs/resources/pipeline-components/sections/resetter_values.yaml | Removes resetter values section docs. |
| docs/docs/resources/pipeline-components/pipeline.yaml | Removes resetter values examples from connector sections. |
| docs/docs/resources/pipeline-components/kafka-source-connector.yaml | Removes resetter values example from source connector doc. |
| docs/docs/resources/pipeline-components/kafka-sink-connector.yaml | Removes resetter values example from sink connector doc. |
| docs/docs/resources/pipeline-components/kafka-connector.yaml | Removes resetter values example from base connector doc. |
| docs/docs/resources/pipeline-components/dependencies/pipeline_component_dependencies.yaml | Removes resetter-values dependency entries. |
| docs/docs/resources/pipeline-components/dependencies/kpops_structure.yaml | Removes resetter fields from documented component field lists. |
| docs/docs/resources/pipeline-components/dependencies/defaults_pipeline_component_dependencies.yaml | Removes resetter-values dependency entry. |
💡 Add Copilot custom instructions for smarter, more guided reviews. Learn how to get started.
|
@copilot in your review, did you read the linked issue? what is the last sentence in the issue? |
Yes — "we can use the new official API instead of the custom streams-bootstrap connector resetter." The |
close #379